Bridging Streaming APIs Behind a REST Proxy
How to expose a simple JSON endpoint when your AI application does not speak REST JSON directly—for example SSE / text/event-stream, WebSockets, or gRPC—so you can integrate it with DynamoEval Custom Applications.
DynamoEval Custom Applications currently integrate with model endpoints over REST and expect a single JSON request/response. Native support for SSE / event-stream, WebSockets, gRPC, and other non-REST protocols is not available yet.
Until that support lands, the recommended approach for anything other than REST JSON is to run a proxy outside of DynamoAI that:
- Accepts DynamoEval (or any REST client) over a normal JSON
POST - Calls your AI application over its native protocol (SSE, WebSocket, gRPC, and so on)
- Aggregates streamed or chunked output into one response when needed
- Returns a single JSON body that DynamoEval can consume
The problem
Many LLM and chat APIs do not return one JSON blob. They may:
- Stream tokens as Server-Sent Events (SSE) /
text/event-stream - Speak WebSockets or gRPC instead of HTTP JSON
POST - Otherwise expose a protocol DynamoEval cannot call directly
That works well for interactive UIs. It does not work well when:
- Evaluation harnesses such as DynamoEval expect a single JSON response
- Sync clients cannot consume event streams or binary RPC protocols
- You need a stable contract even if the upstream protocol changes
Solution: put a REST proxy in front. Clients call a normal POST and get {"content": "..."}. The proxy talks to the upstream protocol, reads every chunk (or RPC response), concatenates the assistant text when streaming, and returns one JSON body.
The reference implementation below focuses on SSE, which is the most common case. The same adapter idea applies to WebSockets and gRPC: accept REST on the outside, speak the native protocol on the inside.
Client ──POST /proxy──► Proxy ──stream POST──► Upstream (SSE)
│
▼
consume stream
aggregate deltas
│
Client ◄──{"content": "..."}──┘
After the proxy is running, register its REST URL as the remote model endpoint in DynamoEval using the standard Custom API Language Models flow (including JSONata transformations if needed). See Custom API Language Model Requirements and Custom Systems Overview for prerequisites, authentication, and transformation basics.
Architecture
| Layer | Role |
|---|---|
| Public REST endpoint | Accepts JSON, returns aggregated JSON |
| Stream forwarder | Opens an HTTP stream to the upstream with httpx |
| Stream consumer | Parses each SSE event and extracts text deltas (application-specific) |
| Error mapper | Maps timeouts and upstream HTTP errors to clear status codes |
Reference stack used in this guide: FastAPI + httpx (async streaming). You can implement the same pattern in any language that supports HTTP streaming.
Authentication
How you handle auth is up to you. Two common patterns:
- Auth at the proxy — Configure DynamoEval’s Custom AI system with credentials for the proxy. The proxy validates those credentials, then authenticates itself to the upstream AI application (for example with a service token it holds).
- Pass-through auth — Configure DynamoEval’s Custom AI system with the AI application’s real credentials. The proxy forwards those headers unchanged to the upstream endpoint.
Either approach works with DynamoEval’s supported auth modes (No Auth, Bearer Token, API Key). Pick the model that matches your security and ops preferences.
Step 1 — Public REST endpoint
Expose a normal POST route. Callers see REST; they never see the stream.
from fastapi import FastAPI, Request, HTTPException
app = FastAPI(title="Streaming Bridge Proxy")
UPSTREAM_URL = "https://your-streaming-api.example.com/v1/chat"
@app.post("/proxy")
async def proxy_chat(request: Request):
try:
return await make_sse_request(UPSTREAM_URL, request)
except HTTPException:
raise
except Exception as e:
return {"content": f"UNCAUGHT_ERROR: {type(e).__name__}: {e}"}
Response shape:
{
"content": "The full assistant reply, assembled from every stream delta."
}
Map this shape to DynamoEval’s expected response format with a JSONata response transformation if needed—for example [content].
Step 2 — Forward the request as a stream
Use httpx.AsyncClient.stream() so the connection stays open while chunks arrive. Do not use a one-shot client.post() that buffers the whole body before you can parse it.
import httpx
from fastapi import Request, HTTPException, Response
TIMEOUT_IN_S = 180.0
async def make_sse_request(url: str, request: Request):
try:
client_body = await request.json()
except Exception:
raise HTTPException(status_code=400, detail="Invalid JSON body received.")
async with httpx.AsyncClient() as client:
try:
async with client.stream(
"POST",
url,
json=client_body,
headers=dict(request.headers),
timeout=TIMEOUT_IN_S,
) as response:
if response.status_code >= 400:
error_body = (await response.aread()).decode("utf-8", errors="replace")
raise HTTPException(status_code=response.status_code, detail=error_body)
content = await consume_stream(response)
return {"content": content}
except httpx.ConnectTimeout:
return Response(status_code=502, content="LLM Connection Timeout")
except httpx.ReadTimeout:
return Response(status_code=504, content="LLM Read Timeout")
except httpx.ConnectError:
return Response(status_code=502, content="LLM Connection Error")
Prefer a long timeout (for example, 180s). Streams finish when generation finishes, not when the first byte arrives.
Step 3 — Consume and aggregate the stream
The consume_stream logic below is an example for one AI application’s SSE event shape. How you parse and aggregate the stream depends entirely on how your AI application streams responses—event names, JSON fields, and framing (data: lines, keep-alives, and so on) will differ. Treat the functions here as a starting point and adapt the extractor to your upstream schema.
SSE typically sends lines like data: {...}. Strip the prefix, parse JSON, and pull out text deltas.
Example event payload (illustrative only):
{"obj": {"type": "message_delta", "content": " Hi"}}
Consumer (example—rewrite extract_delta_content for your API):
import json
import logging
logger = logging.getLogger(__name__)
def extract_delta_content(payload: dict) -> str | None:
"""Pull assistant text from one SSE event. Adapt to your upstream schema."""
obj = payload.get("obj")
if not isinstance(obj, dict):
return None
if obj.get("type") != "message_delta":
return None
content = obj.get("content")
return content if isinstance(content, str) else None
async def consume_stream(response: httpx.Response) -> str:
"""Read SSE lines and concatenate text deltas for this application."""
content_parts: list[str] = []
async for line in response.aiter_lines():
if not line or line.startswith(":"):
continue # keep-alives / comments
if line.startswith("data:"):
line = line[5:].strip()
if not line.strip():
continue
try:
payload = json.loads(line)
except json.JSONDecodeError:
logger.warning("Skipping non-JSON stream line: %r", line)
continue
chunk = extract_delta_content(payload)
if chunk:
content_parts.append(chunk)
return "".join(content_parts)
The aggregation idea is always the same: parse → extract text → join → return once. Only the parsing and field extraction change per application.
Step 4 — Error handling map
| Failure | Typical status | Meaning |
|---|---|---|
| Invalid client JSON | 400 | Bad request body |
Upstream 4xx / 5xx | Same as upstream | Pass through detail when safe |
| Connect timeout | 502 | Couldn’t reach the model service |
| Read timeout | 504 | Stream didn’t finish in time |
| Connect error | 502 | Network / DNS / TLS failure |
Check upstream status before consuming the body so error payloads are not treated as chat deltas.
End-to-end flow
1. Client POSTs JSON to /proxy
2. Proxy validates JSON
3. Proxy opens streamed POST to upstream (headers passed through)
4. If status >= 400 → raise with upstream body
5. Else read SSE events until the stream ends
6. For each event, extract text deltas using your app-specific parser
7. Return {"content": "<joined text>"}
Why this pattern works well
For callers
- One request, one JSON response
- Works with DynamoEval, scripts, and sync HTTP clients
- Stable contract even if upstream streaming format evolves
For operators
- Flexible auth (proxy-owned credentials or pass-through to the AI application)
- Timeouts and connection errors map to clear HTTP codes
- Logging can sit at the proxy without changing clients
Trade-off
- Clients wait until the full answer is generated (no token-by-token UI). That is intentional when the goal is a complete REST response for evaluation, not live streaming.
Minimal dependencies
fastapi>=0.104.0
uvicorn[standard]>=0.24.0
httpx>=0.25.0
Adapting this to your API
- Point
UPSTREAM_URLat your streaming chat endpoint. - Rewrite
extract_delta_content()/consume_stream()to match your SSE event schema (delta,choices[0].delta.content, custom event types, and so on). - Choose an auth pattern (proxy auth vs pass-through) and configure DynamoEval accordingly.
- Keep aggregation and error mapping; the stream consumer is what usually changes.
- Point DynamoEval’s custom model
remote_model_endpointat the proxy’s REST URL and configure auth / JSONata as usual.
Summary
A streaming-to-REST proxy is an adapter:
- Accept a normal JSON POST
- Open an HTTP stream upstream
- Parse and concatenate text deltas (using logic specific to your AI application)
- Return one aggregated JSON body
With FastAPI and httpx’s client.stream(), that bridge is a few focused functions—and it lets DynamoEval (and other tools that only speak REST) talk to backends that only speak SSE.